package com.lhd.app.ods;

import com.lhd.common.utils.MyKafkaUtil;
import com.lhd.common.utils.MysqlUtil;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class TaqiUserApp {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        DataStream<String> ds = MysqlUtil.cdcMysql(env, "gmall", "taqi_user_info");
        ds.print();
        ds.addSink(MyKafkaUtil.getFlinkKafkaProducer("dwd_taqi_user_info"));

        env.execute();
    }
}
